runtime: queue-length limits for the native runtime — scope question - #35
xiaoyu-xyz wants to merge 2 commits into
Conversation
hsliuustc0106
left a comment
There was a problem hiding this comment.
Reviewed commit 1b63ed7431ea4107852600d07799a826773e3bf4.
- [P2] Reserve queue depth before publishing the job. src/frontend/src/worker.rs:217 increments depth after try_send, but the consumer never takes the admission mutex. It can dequeue/process/decrement first, yielding depth=0 in a response and transient usize underflow in health. A standalone CPU probe against the built crate observed four zero-depth answers in 100,000 sequential immediate requests. Reserve before sending and roll back on send failure, or synchronize dequeue/accounting.
- [P2] Start the drain budget when shutdown is signaled. src/frontend/src/main.rs:68 awaits graceful HTTP shutdown before stop_accepting/run_until_drained. In-flight upload/response handlers can consume their entire request budget before the drain timer even begins, and uploads can still submit during this interval. A real binary with DRAIN_MS=50 and TIMEOUT_MS=3000 took 2.917 seconds to exit after SIGTERM with one partial upload. Stop admission on the signal and bound the complete shutdown sequence with the drain deadline.
All 31 ordinary tests plus 1 doctest passed; these probes expose gaps in the existing suite.
1b63ed7 to
3219451
Compare
|
Superseded — see the current status comment. This comment is kept only so the thread does not lose a reply that was read. Two of its claims did not survive re-checking, and I am not leaving them standing as if they had:
Nothing else in it was withdrawn — the two P2 findings were fixed as described, and the shutdown and capacity findings that came afterwards are covered in the current comment. |
Levius-Fubuki
left a comment
There was a problem hiding this comment.
Reviewed 3219451e57ab7ad96fc92ae47d4636fe404f9c7b, including the incremental changes on #34 and the two previous findings.
The depth reservation now precedes publication and rolls back on both failed-send paths, addressing the original counter-underflow mechanism. However, the shutdown follow-up introduces a reproducible service-lifetime regression. Changes are requested for the three issues below.
Formatting passed. Clippy, workspace tests and release build passed on Linux/Rust 1.98.1 after rebuilding in a dedicated target directory (an initial shared-target run had stale cross-checkout artifacts and is not treated as a PR failure). Independent probes reproduce the no-signal exit and invalid-mode fallback on Linux and macOS; a gated CPU worker reproduces the capacity breach. No GPU is needed for these cases. This remains a draft and also needs its existing merge conflicts resolved.
| // Responses in flight get the drain budget to finish. A handler still stuck on a slow | ||
| // upload at the deadline is not going to finish, and holding the process open for it | ||
| // would mean the budget bounded nothing. | ||
| let deadline = Instant::now() + running.drain_budget(); |
There was a problem hiding this comment.
[P1] Start the drain deadline on the shutdown signal, not at startup
This deadline is created as soon as the listener starts and wraps the entire engine::serve future. Consequently every native process exits after the drain duration even if no SIGTERM/SIGINT is sent (10 seconds by default). With OMNI_SYSTEMONE_ENGINE=passthrough, OMNI_SYSTEMONE_DRAIN_MS=200 and OMNI_JEV_BIND=127.0.0.1:0, the release binary exited 0 after 0.206 seconds on Linux, logging shutdown exceeded its drain budget; the probe sent no signal. Wait for the signal while serving normally, then close admission and compute one absolute deadline for graceful HTTP shutdown plus worker drain. Add a test that remains healthy beyond the drain duration before sending any signal.
| // answered and released — counting afterwards would release a slot that had not | ||
| // been taken yet, which reads as a wrong depth and wraps the counter. The window is | ||
| // small but real, and `depth_is_reserved_before_the_job_is_published` caught it. | ||
| self.depth.fetch_add(1, Ordering::AcqRel); |
There was a problem hiding this comment.
[P2] Enforce the documented limit on all outstanding requests
This increment reserves accounting but never checks the configured outstanding-request limit; only mpsc::try_send enforces capacity. Once the worker dequeues a job, the channel slot is free even while inference still runs. A deterministic probe with capacity 1 holds the first inference behind a gate, waits until it has started, then submits a second job: it is accepted and report().depth becomes 2. Both Options::queue_capacity and the README promise refusal once that many accepted requests are unfinished. Check/reserve the outstanding count under the admission lock (with rollback on send failure), and cover a running-plus-queued case instead of relying on whether the consumer has happened to dequeue.
|
|
||
| async fn run() -> Result<(), BoxError> { | ||
| let bind = setting("OMNI_JEV_BIND", Config::DEFAULT_BIND.to_owned())?; | ||
| if env::var("OMNI_SYSTEMONE_ENGINE").as_deref() != Ok("passthrough") { |
There was a problem hiding this comment.
[P2] Reject unsupported engine names instead of silently selecting forwarding
Every value other than passthrough enters forwarding mode here. With OMNI_SYSTEMONE_ENGINE=typo-engine, the built binary starts and logs forwarding to http://127.0.0.1:8000/, although the README says unknown modes are startup errors. A configuration typo can therefore route inference to a different backend instead of failing visibly. Match supported values explicitly, preserve the documented unset/default behavior, and return a configuration error for unknown values.
3219451 to
a9f7b83
Compare
|
Rechecked |
a0e86f4 to
c88dc14
Compare
RFC ThinkFlowLab#14 puts the request lifecycle in src/frontend/ and the model behind "a small engine interface", and #3's acceptance criteria depend on that interface existing. Neither half is on main today. This adds the interface and the HTTP service behind it: the Engine trait, its readiness/answer/error types, and engine::app serving /v1/systemone and /health with a request budget measured from the headers, a body limit, and one status code per failure mode. Nothing parses the decision envelope, so a field this crate has never heard of survives and a model's error text cannot produce invalid JSON. submit takes owned bytes and returns a channel, which is the only shape that works for a worker thread with a queue behind it; the reply carries the queue depth to publish, so x-queue-depth is measured by the engine rather than guessed by the transport. No model, no GPU, and no change to main.rs or the forwarding path: the thread, the bounded queue and the drain that plug into this boundary come next.
c88dc14 to
63b641c
Compare
|
@Levius-Fubuki — the shutdown case you flagged is fixed, and you were right about why my test missed it. Your review was against The measurement you reported reproduces here on Why the wait cannot be signalled away. I instrumented the two halves separately: the shutdown future finished at 1.01 s, and My test was the reason I did not see it. It hung up before signalling, and the bug only appears while the connection is open — I wrote a test for a bug and then arranged the precondition away. It now holds Measured with the connection held open through the whole signal-to-exit window:
The bound now runs from the signal and only from the signal: the shutdown future closes admission, starts the budget, drains the worker, and serves a bounded grace period for responses already in flight; the outer wait ends on whichever comes first — Two wrong attempts before this one, both of which looked right: timing the serve future from startup re-introduced the no-signal exit, and a Also rebased onto current |
|
Rechecked |
63b641c to
9876a7b
Compare
|
@Levius-Fubuki — one absolute deadline now, and the regression test is tight enough to stay red on the additive version. What changed. The deadline is taken once, at the signal, and every step after it is bounded by that instant rather than by a fresh period:
That also fixed a regression I introduced while bounding this: because the window was the whole budget, a binary with A regression test for the overrun itself. One thing I have to be straight about. I cannot reproduce the numbers you measured on What I can state as measured, with the connection held open through the shutdown:
And the structural claim, which is what your instruction asked for and what the source now shows: there is one deadline, taken at the signal; the worker drains to it; the transport's window is clamped to it; nothing adds a period measured afterwards. If your probe still shows an overrun on Verified on a clean checkout of |
|
Superseded by #34 (reopened with the placement revision; #91 was a duplicate and is closed). Closing this branch. The placement finding on #34 applies here more directly than anywhere: this is the worker-thread queue and lifecycle that the finding is about. The contract it would implement now lives in A correction to what I first wrote here. I said the tests worth keeping from this branch are in #34. That was too broad, and I checked rather than leave it standing:
So nine of those thirteen are not in #34; they belong to the queue and lifecycle, and they went out with it. Two of the seven rows above were also findings from review that were fixed here first — the capacity definition and the depth reservation — and the second has a known limit that has not changed: the underflow interleaving was never reproduced by my own test, only by the reviewer's probe. What this branch leaves behind, then, is a specification of the behaviour a runtime queue would need: readiness that follows warmup, a bound that never counts a slot it did not take, work that is dropped rather than started when its caller has gone or its budget expired, a failure that retires the engine instead of serving from unknown state, and a drain that reports when it could not finish. If queue-length limits are picked up, those are the cases to cover, and the branch is here to read. |
|
The question this PR is open for, in one place, so it can be answered without reading the stale diff.
Nothing is pending on an answer — this branch's head is stale and would be rewritten against |
|
Rechecked |
9876a7b to
cbbeedf
Compare
WARNING: this head is stale and is not a merge candidate. It implements the queue in the frontend, which review asked to move to the worker-side runtime, and it predates the contract it would build on (`Engine` now lives in src/runtime). It is kept because the queue-length question in the PR body needs an answer before the work is redone against current `main`, and because the tests here specify behaviour a runtime queue would have to satisfy. What it is, for reading: - worker.rs: one thread owns the engine, a bounded queue owns admission, and readiness, cancellation, expiry and drain are the transport's. A slot is taken before a job is handed over and rolled back if the handover fails; a job whose caller has gone or whose budget expired is dropped instead of started; a panic in model code answers its caller and retires the engine. - passthrough.rs: a stand-in engine that echoes requests, so the path can run with no model and no GPU. - worker_lifecycle.rs: the cases above, each with a failure mode behind it. The last change synchronizes three depth assertions against the worker's own reply ordering. The worker sends a reply before it releases the slot that reply held, so reading depth immediately after `reply` could still see the slot in use and fail about once in a full workspace run. Depth only falls as the worker finishes, so the checks now wait for zero with a deadline. I could not reproduce the window locally in twelve instrumented runs, so the reason for the change is the ordering in the code and the reported failure, not my own observation.
cbbeedf to
c423100
Compare
|
Synchronized on The three depth checks now wait for zero rather than sampling once. Depth only falls as the worker finishes, so waiting is what checks the property actually being claimed — every accepted request releases its slot — and a deadline still fails the test if some request never does. I converted the same shape in the capacity test and the storm test while I was there. One thing I have to be straight about: I could not reproduce the window locally. I instrumented the read to print the depth at the moment the reply arrived and ran it twelve times; it read zero every time. So the reason for the change is the ordering in the source — After the change: 40 isolated runs of that test, 15 runs of the whole suite, and 12 full workspace runs, all passing (57 passed, 9 ignored). On the stale-implementation note: agreed, and the branch says so now — the commit message opens by stating it is not a merge candidate and why, and the PR body asks the scope question rather than proposing the diff. No work is pending on it. |
Purpose
Queue-length limits for the native runtime — whether they are wanted, and whether now is the time.
omni-runtimeowns admission asSerialScheduler: FIFO, one execution unit at a time, no queue-length limit. Its README lists this explicitly:This PR is that planned item, in the layer the architecture assigns it to. The contract it would build on landed in #34:
Engineinsrc/runtime, withReport { depth, capacity, rejected }shaped to describe exactly this.What #34 already has, and what would be left here
src/runtime/src/engine.rstests/runtime/engine.rs503vs504linetests/frontend/in_process.rscapacityrequests are outstanding, running or queuedspawnreturning only after load and warmup, so readiness means somethingThe behaviour being proposed
Specified from what that implementation had to get right, and from two review findings it fixed:
usizethat wraps, and/healththen reports a number nobody can act on.spawnreturns only once the model has loaded and warmed up, so a handle cannot observe a half-loaded engine.The question
omni-runtimenow, or should the limits stay planned until a worker needs them? Two native workers currently run with concurrency one and no limit, and neither has reported back-pressure as a problem.capacityas outstanding-request count with refusal at the bound,Reportpublishing depth/capacity/rejected, readiness tied to warmup, failure retiring the engine? Anything you would draw differently is cheaper to change before a model depends on it.SerialScheduler, or a separate admission type beside it?SerialScheduleris 40 lines and knows nothing about bodies or deadlines; a bounded queue that carries a deadline and reports depth is a different thing.No work is pending on this: the current head is stale and would be rewritten against
mainand #34 if the answers say to do it. Asking first so the rewrite lands once.